package com.edu.realtime;

import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;

import com.edu.realtime.util.MyKafkaUtil;

/**
 * Created on 2022/10/18.
 *
 * @author Topus
 */
public class KafkaConsumerDemo {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(4);
        FlinkKafkaConsumer<String> kafkaConsumer = MyKafkaUtil.getKafkaConsumer("topic_log", "topic_group");

        DataStreamSource<String> source = env.addSource(kafkaConsumer);
        source.print();

        env.execute();
    }

}
